package com.mjf.day2

import org.apache.flink.streaming.api.scala._

/**
 * 从自定义数据源消费数据
 */
object ConsumerFromSensorSource {
  def main(args: Array[String]): Unit = {

    val env: StreamExecutionEnvironment = StreamExecutionEnvironment.getExecutionEnvironment
    env.setParallelism(1)

    // 调用addSource方法
    val stream: DataStream[SensorReading] = env.addSource(new SensorSource)

    stream.print()

    env.execute("ConsumerFromSensorSource")

  }
}
